FlowStreamType
FlowStream
FlowStream<'env, 'error, 'value> describes a cold, pull-based sequence of effectful values. It uses the same
environment, typed failure channel, cancellation, and runtime scopes as Flow.
Choose FlowStream when values should be handled incrementally instead of collected before work begins. Typical
sources include process output, paginated APIs, subscriptions, and host-specific file or network adapters.
A stream does nothing until a terminal operation such as runFold,
runForEachFlow, or
runCollect produces a Flow and that Flow is run. Pulling provides backpressure: downstream requests each next value, so upstream cannot run ahead
without an operator explicitly introducing bounded concurrency. The time-based operators (groupedWithin, throttle,
debounce, switchMapFlow) and buffer are such operators: they read ahead into a queue so they can react to
time or to a newer value.
A stream that reads a file, socket, or cursor should own it with using, which releases
the resource however consumption ends. To take a single result without draining the stream, use
runTryHead.
Streams connect to the concurrency types in both directions. fromDequeue reads a
queue or hub subscription until it is shut down and drained, and fromHub subscribes for
the life of the stream. runIntoQueue and runIntoHub
feed a stream into them. mergePar runs several streams concurrently, and
fromSchedule turns a schedule into a tick stream. See the Queue
and Hub guides.
Start here
The Streams guide introduces the model. Getting started builds a
complete pipeline, while Batching and parallelism explains the scheduling
difference between strict batches and mapFlowPar.
Use incremental terminal operations for large or unbounded streams. runCollect deliberately retains every emitted
value and should be reserved for known finite inputs.
Specification
Summary
| Name | Signature | Synopsis |
|---|---|---|
| Union cases | ||
| FlowStream | FlowStream 'env -> CancellationToken -> Execution | No description available. |
| Operations | ||
| localEnv | FlowStream.localEnv mapping (FlowStream operation) | Runs a stream with an environment derived from its caller's environment. |
| unfoldFlow | FlowStream.unfoldFlow step initialState | Creates a cold stream by repeatedly running an effectful state transition. |
| fromDequeue | FlowStream.fromDequeue queue | Creates a stream that takes values from a queue or hub subscription until it is shut down and drained. |
| fromHub | FlowStream.fromHub strategy hub | Creates a stream of the values published to a hub, subscribing for the life of the stream. |
| fromSchedule | FlowStream.fromSchedule schedule | Creates a stream that emits a schedule's outputs, each after the delay the schedule chooses. |
| runIntoQueue | FlowStream.runIntoQueue queue stream | Runs a stream and offers every value to a queue, waiting while a bounded queue is full. |
| runIntoHub | FlowStream.runIntoHub hub stream | Runs a stream and publishes every value to a hub. |
| buffer | FlowStream.buffer strategy stream | Runs the upstream ahead of its consumer in its own fiber, buffering values with strategy. |
| mergePar | FlowStream.mergePar streams | Runs several streams concurrently and emits their values as they arrive. |
| fromSeq | FlowStream.fromSeq values | Creates a stream from a synchronous sequence of values. |
| empty | FlowStream.empty | Creates an empty stream. |
| singleton | FlowStream.singleton value | Creates a stream containing one value. |
| fromFlow | FlowStream.fromFlow flow | Creates a one-element stream from an effectful value. |
| runForEach | FlowStream.runForEach action (FlowStream op) | Executes the stream and performs a synchronous action for each successful value. |
| map | FlowStream.map f stream | Transforms the successful values of a stream using the provided function. |
| mapError | FlowStream.mapError mapper stream | Transforms the typed error channel of a stream. |
| filter | FlowStream.filter predicate stream | Keeps values that satisfy a predicate. |
| choose | FlowStream.choose chooser stream | Maps and filters values in one operation. |
| tapFlow | FlowStream.tapFlow action stream | Runs an effect for each value before emitting the original value. |
| mapFlow | FlowStream.mapFlow mapper stream | Transforms every value with a Flow effect. |
| mapFlowPar | FlowStream.mapFlowPar parallelism mapper stream | Maps values with a continuously replenished, bounded set of child fibers. |
| mapFlowParUsing | FlowStream.mapFlowParUsing parallelism resource mapper stream | Like mapFlowPar, giving each running mapping a resource from a pool that holds at most parallelism of them. |
| groupedWithin | FlowStream.groupedWithin size window stream | Groups values into lists of at most size, emitting early when window passes. |
| debounce | FlowStream.debounce quiet stream | Emits a value only once quiet passes without a newer one. |
| throttle | FlowStream.throttle interval stream | Emits at most one value per interval, keeping the latest value when faster. |
| switchMapFlow | FlowStream.switchMapFlow mapper stream | Maps each value to a flow, interrupting the previous flow when a newer value arrives. |
| chunkBySize | FlowStream.chunkBySize size stream | Groups consecutive values into non-empty lists of at most size elements. |
| take | FlowStream.take count stream | Emits at most count values. |
| skip | FlowStream.skip count stream | Skips the first count values. |
| takeWhile | FlowStream.takeWhile predicate stream | Emits values while a predicate remains true. |
| skipWhile | FlowStream.skipWhile predicate stream | Skips values while a predicate remains true. |
| indexed | FlowStream.indexed stream | Emits each value paired with its zero-based index. |
| scan | FlowStream.scan folder initial stream | Emits successive accumulator states. |
| distinctUntilChangedBy | FlowStream.distinctUntilChangedBy projection stream | Suppresses consecutive duplicate values according to a projection. |
| append | FlowStream.append (FlowStream right) (FlowStream left) | Concatenates two streams, evaluating the second only after the first ends. |
| collect | FlowStream.collect mapper stream | Maps each value to a stream and concatenates the resulting streams. |
| zip | FlowStream.zip rightStream leftStream | Pairs values from two streams until either stream ends. |
| runFold | FlowStream.runFold folder initial stream | Folds a stream into one value inside Flow. |
| runCollect | FlowStream.runCollect stream | Collects all emitted values into a list. |
| runDrain | FlowStream.runDrain stream | Consumes a stream and ignores its values. |
| runForEachFlow | FlowStream.runForEachFlow action stream | Runs an effectful action for every stream value. |
| repeatFlow | FlowStream.repeatFlow flow | Creates a stream that emits the result of running flow again for every pull, forever. |
| using | FlowStream.using resource build | Creates a stream that owns a resource: acquired when consumption starts, released when it ends. |
| runTryHead | FlowStream.runTryHead stream | Returns the first value, or None for an empty stream, then stops the stream. |
| runTryLast | FlowStream.runTryLast stream | Consumes the whole stream and returns its last value, or None for an empty stream. |
| runCount | FlowStream.runCount stream | Consumes the whole stream and returns how many values it emitted. |
Union cases
Operations
Parameters
| Name | Type | Description |
|---|---|---|
| mapping | 'outerEnvironment -> 'innerEnvironment | |
| FlowStream operation | FlowStream<'innerEnvironment, 'error, 'value> |
Returns
Parameters
| Name | Type | Description |
|---|---|---|
| step | 'state -> Flow<'env, 'error, ('value * 'state) option> | Returns Some(value, nextState) or None when the stream is complete. |
| initialState | 'state | The state used for the first pull. |
Returns
Verification Examples
FlowStream.unfoldFlow (fun n -> Flow.ok (if n < 3 then Some(n, n + 1) else None)) 0Parameters
| Name | Type | Description |
|---|---|---|
| queue | Dequeue<'value> |
Returns
Verification Examples
flow {
let! (jobs: Queue<string>) = Queue.bounded 8
do! jobs |> Queue.offerAll [ "a"; "b" ] |> Flow.ignore
do! Dequeue.shutdown jobs
return! jobs |> FlowStream.fromDequeue |> FlowStream.runCollect
}Parameters
| Name | Type | Description |
|---|---|---|
| strategy | QueueStrategy | |
| hub | Hub<'value> |
Returns
Verification Examples
readings
|> FlowStream.fromHub (QueueStrategy.Sliding 1)
|> FlowStream.runForEach displayParameters
| Name | Type | Description |
|---|---|---|
| schedule | Schedule<'env, unit, 'output> |
Returns
Verification Examples
// A tick every 50 ms, aligned to when the stream started.
Schedule.fixedRate (TimeSpan.FromMilliseconds 50.0) |> FlowStream.fromScheduleParameters
| Name | Type | Description |
|---|---|---|
| queue | Queue<'value> | |
| stream | FlowStream<'env, 'error, 'value> |
Returns
Verification Examples
deviceReadings |> FlowStream.runIntoQueue controlInputsstrategy.Parameters
| Name | Type | Description |
|---|---|---|
| strategy | QueueStrategy | |
| stream | FlowStream<'env, 'error, 'value> |
Returns
Verification Examples
samples |> FlowStream.buffer (QueueStrategy.Sliding 1) |> FlowStream.runForEachFlow renderParameters
| Name | Type | Description |
|---|---|---|
| streams | FlowStream<'env, 'error, 'value> list |
Returns
Verification Examples
[ readerA; readerB; readerC ] |> FlowStream.mergePar |> FlowStream.runIntoQueue controlInputsParameters
| Name | Type | Description |
|---|---|---|
| values | 'value seq | The sequence of values to be emitted by the stream. |
Returns
Verification Examples
FlowStream.fromSeq [1..10]
|> FlowStream.runCollect
|> Flow.run ()Parameters
| Name | Type | Description |
|---|---|---|
| action | 'value -> unit | The function to execute for each value emitted by the stream. |
| FlowStream op | FlowStream<'env, 'error, 'value> |
Returns
Verification Examples
FlowStream.fromSeq ["a"; "b"; "c"]
|> FlowStream.runForEach (printfn "%s")
|> Flow.run ()Parameters
| Name | Type | Description |
|---|---|---|
| f | 'v -> 'w | The function to transform each value. |
| stream | FlowStream<'env, 'error, 'v> | The stream whose values should be transformed. |
Returns
Verification Examples
let stream = FlowStream.fromSeq [1; 2; 3] |> FlowStream.map (fun n -> n * 2)mapFlowPar, giving each running mapping a resource from a pool that holds at most parallelism of them.Parameters
| Name | Type | Description |
|---|---|---|
| parallelism | Parallelism | |
| resource | Resource<'env, 'error, 'resource> | |
| mapper | 'resource -> 'value -> Flow<'env, 'error, 'next> | |
| stream | FlowStream<'env, 'error, 'value> |
Returns
Verification Examples
commits |> FlowStream.mapFlowParUsing (Parallelism.bounded 4) openReader (fun reader commit -> diff reader commit)size, emitting early when window passes.Parameters
| Name | Type | Description |
|---|---|---|
| size | int | |
| window | TimeSpan | |
| stream | FlowStream<'env, 'error, 'value> |
Returns
Verification Examples
events |> FlowStream.groupedWithin 100 (TimeSpan.FromSeconds 1.0)quiet passes without a newer one.Parameters
| Name | Type | Description |
|---|---|---|
| quiet | TimeSpan | |
| stream | FlowStream<'env, 'error, 'value> |
Returns
Verification Examples
keystrokes |> FlowStream.debounce (TimeSpan.FromMilliseconds 300.0)interval, keeping the latest value when faster.Parameters
| Name | Type | Description |
|---|---|---|
| interval | TimeSpan | |
| stream | FlowStream<'env, 'error, 'value> |
Returns
Verification Examples
progress |> FlowStream.throttle (TimeSpan.FromMilliseconds 100.0)Parameters
| Name | Type | Description |
|---|---|---|
| mapper | 'value -> Flow<'env, 'error, 'next> | |
| stream | FlowStream<'env, 'error, 'value> |
Returns
Verification Examples
queries |> FlowStream.debounce (TimeSpan.FromMilliseconds 200.0) |> FlowStream.switchMapFlow searchParameters
| Name | Type | Description |
|---|---|---|
| projection | 'a -> 'b | |
| stream | FlowStream<'c, 'd, 'a> |
Returns
Verification Examples
stream |> FlowStream.distinctUntilChangedBy idParameters
| Name | Type | Description |
|---|---|---|
| FlowStream right | FlowStream<'a, 'b, 'c> | |
| FlowStream left | FlowStream<'a, 'b, 'c> |
Returns
Verification Examples
first |> FlowStream.append secondParameters
| Name | Type | Description |
|---|---|---|
| mapper | 'a -> FlowStream<'b, 'c, 'd> | |
| stream | FlowStream<'b, 'c, 'a> |
Returns
Verification Examples
stream |> FlowStream.collect FlowStream.fromSeqParameters
| Name | Type | Description |
|---|---|---|
| action | 'value -> Flow<'env, 'error, unit> | |
| stream | FlowStream<'env, 'error, 'value> |
Returns
Verification Examples
stream |> FlowStream.runForEachFlow saveflow again for every pull, forever.Parameters
| Name | Type | Description |
|---|---|---|
| flow | Flow<'env, 'error, 'value> |
Returns
Verification Examples
FlowStream.repeatFlow readSensor |> FlowStream.takeWhile (fun reading -> reading.Ok)Parameters
| Name | Type | Description |
|---|---|---|
| resource | Resource<'env, 'error, 'resource> | The resource to acquire. |
| build | 'resource -> FlowStream<'env, 'error, 'value> | Builds the stream from the acquired resource. |
Returns
Verification Examples
let lines path =
FlowStream.using (Resource.create (Flow.fromBlocking (fun _ -> File.OpenText path)) (fun reader _ -> reader.Dispose(); Task.CompletedTask))
(fun reader -> FlowStream.repeatFlow (Flow.fromBlocking (fun _ -> reader.ReadLine())) |> FlowStream.takeWhile (isNull >> not))
